package com.atlocal.app;

import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject;
import com.atlocal.base.impl.FlinkBaseApiImpl;
import com.atlocal.utils.MyKafkaUtil;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.functions.ProcessFunction;
import org.apache.flink.streaming.api.functions.sink.SinkFunction;
import org.apache.flink.util.OutputTag;

/**
 * @ClassName ProDataToTest
 * @Description TOO
 * @Author kongjiangjiang
 * @Date 2024-01-02 15:45
 * @Version 1.0
 **/
public class ProDataToTest extends FlinkBaseApiImpl {
    public static void main(String[] args) throws Exception {
        //动态读取生产-kafka数据
        DataStreamSource<String> kafkaStream = MyKafkaUtil.getKafkaSourceDy(env, "test", true);


        OutputTag<String> pro_dsg2kfk_zssys_web_ply_base = new OutputTag<String>("pro_dsg2kfk_zssys_web_ply_base") {
        };



        env.execute();
    }
}
